D16 講完資料怎麼在 DataFusion 內部流:拉、一批一批、跑在 tokio 上,pipeline breaker 是 async 的等待點。今天緊接著問一個沒答案就不能上線的問題:
同一個 tokio runtime 上跑好幾個 pipeline breaker,它們都想要記憶體,記憶體只有一份,誰該讓?
這問題乍看是資源管理小事,但答錯了會直接爆整個 executor process。DataFusion 給的答案是一個叫 MemoryPool 的 trait,兩種預設實作代表兩種公平觀。今天把它拆開看,順便回扣 Comet 一個真痛點:Spark 那邊有 UnifiedMemoryManager 在管 JVM 記憶體、Rust 這邊又有一個 pool,兩個池怎麼對齊。
假設要跑這一句:
SELECT l_orderkey, SUM(l_extendedprice)
FROM lineitem JOIN orders ON l_orderkey = o_orderkey
WHERE o_orderdate > '1998-01-01'
GROUP BY l_orderkey
ORDER BY 2 DESC
跑一次 EXPLAIN,會看到大概長這樣的樹:
SortExec
└── AggregateExec (Final)
└── HashJoinExec
├── build side: DataSourceExec (orders)
└── probe side: DataSourceExec (lineitem)
D16 講過三種是 pipeline breaker——下游要等上游把資料物化完才能開始拉。這裡剛好三個都有:
orders 的 hash table 要建好,lineitem 才能開始 probe三個大戶同時要記憶體。整台機器 heap 只有那麼多,如果都不管,OS 一個 OOM 直接把 process 幹掉——連錯誤訊息都沒有,只在 syslog 留一行 killed。
所以問題不是「記憶體不夠怎麼辦」,而是不夠的時候誰該讓、讓多少、由誰來判斷。
寫過 C / Rust 的人第一直覺會是「不是有 allocator 嗎,malloc / new 出錯就處理啊」。可是 allocator 只知道「有沒有那麼多 bytes 給你」,不知道「這個 700 MB 是 SortExec 想繼續讀,還是 HashJoin 想蓋 hash table」,也不知道「誰有 spill 路徑、誰沒有」。
要協調,只能有一個更高一層的東西集中看——這就是 MemoryPool:
MemoryConsumer
reserve(size) 拿一個 MemoryReservation
grow(additional),吃不完 shrink(released)
ResourcesExhausted error這裡有個很重要的分工:pool 不會自己去 spill 別人。它只做「批准或拒絕」,spill 是 consumer 自己的事。理由是 pool 不知道別人的 state 長怎樣、能不能 spill、spill 到哪裡——這些只有 consumer 自己清楚。
用一個小類比:pool 像公司財務,運算子像各部門。財務只管「本月預算超了沒」,超了就打回申請單。部門收到打回會自己決定要延後案子、還是找別的方案,財務不會幫你決定要砍哪個案子。
這個「批准 / 拒絕」跟「怎麼處理拒絕」分開的設計,是後面 Greedy / Fair 兩種策略能長出來的前提。
DataFusion 內建兩個主要實作:GreedyMemoryPool 跟 FairSpillPool。名字取得很直白,但取捨要拆開講。
規則超簡單:加總所有 reservation,超過 limit 就拒絕新的。
先跑的運算子想拿多少拿多少,後跑的一發現額度用完就得馬上 spill。Sort 已經吃了 700 MB、HashJoin 想 grow 200 MB 卻只剩 100 MB 額度——直接拒絕,HashJoin 去 spill,Sort 完全不受影響。
好處是簡單、可預測,適合單 query 或只有一個大戶運算子的場景。壞處是不公平:Sort 明明有 spill 路徑,卻因為它先來就佔盡優勢,HashJoin 明明剛開始蓋 hash table,spill 就會導致後面的 probe 得從 disk 讀。
規則稍微複雜:把「可 spill 的部分」平均分。
關鍵點:Fair 不是「先 spill 誰誰吃虧」,是「已經吃最多的先讓」。
回到剛才的例子。假設 pool 上限 1 GB、Sort 已用 700 MB、HashJoin 用 200 MB。新的 grow 進來時:
這個切法讓「有 spill 路徑」變成一種義務,而不是運氣——你有 spill 就要有讓的準備,不是你先來就可以霸佔。
沒有絕對答案。單 query 的批次任務用 Greedy 反而穩定,因為 spill 越少越好。多個 spillable 運算子並存(例如同時有 sort + hash join build)Fair 比較合理,不會出現先來的把後來的擠到磁碟。
DataFusion 預設是 GreedyMemoryPool,Comet 那邊會覆蓋這個選擇,等下講。
DataFusion 還有第三個工具,叫 TrackConsumersPool——它不是獨立的策略,是包在 Greedy 或 Fair 外面的裝飾器。
功能是:pool 拒絕 reservation 時,錯誤訊息會附上「目前用最多的前 N 個 consumer 各自吃了多少」。這個很實用,OOM 錯誤如果只寫「記憶體滿了」,你根本不知道找誰——是 sort 太貪、還是 hash join 建了太大的表、還是有個 UDF 洩漏。
這個 pool 的近似點在:只維護 top-N 名單,小額配置不追。理由是實務上前幾名通常吃掉八成,追全部要多維護一份 map、每次 reserve / release 都要更新,成本不划算。
問題就出在這個「通常」——不是「總是」。
思考題:什麼場景會讓「只追大戶」這個近似出事?
一個提示:想像一條 query 用了 200 個 filter,每個 filter 各自 buffer 一小塊 batch,單個都排不進 top-10,加起來卻吃掉 300 MB。OOM 時錯誤訊息只會指到那個吃 500 MB 的 sort,但真正該優化的可能是那 200 個小 filter 為什麼要各留一份。
抓錯兇手不會直接害你當機,但會害你花一整個下午在錯的地方調 config。這是近似的隱性成本,發生率不高但一發生就很久。
把上面拼起來,一條 SQL 的 memory 管理路徑:
MemoryPool (limit = 1 GB, Fair 策略)
│
┌──────────┼──────────────────────┐
│ │ │
register register register
│ │ │
▼ ▼ ▼
SortExec HashJoinExec AggregateExec
(spillable) build side (spillable)
(spillable)
[執行到一半] [執行到一半]
grow(200 MB) grow(300 MB)
│ │
▼ ▼
┌──────────────────────────────────────┐
│ pool: 我看看 │
│ Sort 已 700 MB (70%) │
│ HashJoin 已 200 MB (20%) │
│ AggregateExec 已 100 MB (10%) │
│ Fair 份額 = 1000 / 3 ≈ 333 MB │
│ Sort 超額最多 → 新 grow 拒絕 │
│ (或選 Sort 那邊拒絕 → 逼 Sort spill) │
└──────────────────────────────────────┘
│
▼
ResourcesExhausted → SortExec 呼叫自己的 spill 路徑
→ 寫一個 run 到 disk → 釋放 reservation → pool 額度回來
幾個實務細節:
D15 接點四留了一個問題——JVM 那邊的記憶體歸 Spark 管,Rust 這邊的記憶體歸 DataFusion 管,兩邊怎麼對齊。今天可以正式交代。
Spark 這邊有一個 UnifiedMemoryManager:execution(sort、hash、agg)跟 storage(cache、broadcast)共享一塊 JVM heap,兩塊會動態借用。這是 2015 年 Project Tungsten 定下來的模型,沿用至今。
Rust 這邊有 DataFusion 的 MemoryPool。Comet 每一個 native operator 都跑在這個 pool 底下。
兩者不對齊會出這幾種事:
Comet 的處理方向是讓 Rust 側的 MemoryPool 變成 Spark manager 的下級:
這條路的代價是每個 reservation 都 pay 一次 JNI 呼叫。D16 講 block_on 是每個 batch 一次來回,那個頻率已經不低;MemoryPool 的 reservation 可能一個 batch 內 grow / shrink 好幾次——JNI 頻率更高。
Comet 從 0.11 開始花力氣減少這個往返(例如 batch reservation、預先撥額度),但結構性成本不會消失,只能攤平。這件事會在 D26 深挖 JNI 邊界時再回來看細節。
MemoryPool 是 DataFusion 把「記憶體協調」抽出來的一層:不做 spill 決策,只做批准 / 拒絕,spill 的怎麼辦交給 consumer。Greedy 跟 Fair 兩種策略代表兩種公平觀,TrackConsumersPool 是包在外面的追蹤器、只留大戶名單當近似。Comet 這邊要把 native pool 對齊到 Spark 的 UnifiedMemoryManager,代價是每個 reservation 都要跨 JNI 問一次。
明天 D18 接執行細節的第三題:Join 演算法怎麼選?hash、merge、symmetric hash、nested loop 各自的取捨在哪,planner 憑什麼在這幾個裡挑一個?其中一個很硬的因素就是今天的主題——這個演算法要多少記憶體、有沒有 spill 路徑。
那就明天見~
MemoryPool trait 文件: https://docs.rs/datafusion/latest/datafusion/execution/memory_pool/trait.MemoryPool.html
GreedyMemoryPool API 文件: https://docs.rs/datafusion/latest/datafusion/execution/memory_pool/struct.GreedyMemoryPool.html
FairSpillPool API 文件: https://docs.rs/datafusion/latest/datafusion/execution/memory_pool/struct.FairSpillPool.html
TrackConsumersPool API 文件: https://docs.rs/datafusion/latest/datafusion/execution/memory_pool/struct.TrackConsumersPool.html